{T}

mysql2 原生驱动

mysql2 是一个为 Node.js 设计的现代化 MySQL 客户端库,以其卓越的性能和开发者友好性而著称。作为广受欢迎的 mysql 库的直接替代品和精神续作,它提供了显著的速度提升和一系列高级功能。

系统架构

架构图

code
┌─────────────────────────────────────────────────────────────┐
│                        Node.js 应用                          │
└────────────────────────┬────────────────────────────────────┘
                         │
                         ▼
┌─────────────────────────────────────────────────────────────┐
│                      mysql2 驱动层                            │
│  ┌──────────────────────────────────────────────────────┐  │
│  │              Connection Pool (连接池)                 │  │
│  │  ┌──────────┐ ┌──────────┐ ┌──────────┐            │  │
│  │  │Connection│ │Connection│ │Connection│  ...       │  │
│  │  │    1     │ │    2     │ │    3     │            │  │
│  │  └──────────┘ └──────────┘ └──────────┘            │  │
│  └──────────────────────────────────────────────────────┘  │
│  ┌──────────────────────────────────────────────────────┐  │
│  │          Protocol Parser (协议解析器)                 │  │
│  │  • MySQL 协议编解码                                   │  │
│  │  • 预处理语句管理                                     │  │
│  │  • 结果集流式处理                                     │  │
│  └──────────────────────────────────────────────────────┘  │
└────────────────────────┬────────────────────────────────────┘
                         │
                         ▼
┌─────────────────────────────────────────────────────────────┐
│                    MySQL 数据库服务器                        │
│  • 数据存储与管理                                            │
│  • 查询解析与执行                                            │
│  • 事务处理                                                  │
└─────────────────────────────────────────────────────────────┘

核心模块

模块说明主要功能
连接池 (Connection Pool)管理数据库连接的生命周期连接复用、负载均衡、连接泄漏检测
连接 (Connection)单个数据库连接实例执行查询、事务管理
预处理语句 (Prepared Statement)SQL 模板与参数分离防止 SQL 注入、提升性能
流式查询 (Stream Query)大数据集分批处理降低内存消耗、实时数据处理
Promise 支持异步操作封装async/await 语法支持

核心特性

1. 卓越的性能

mysql2 的协议解析器是从头开始编写的,其性能远超原始的 mysql 库:

  • 快速解析:优化的二进制协议解析器
  • 批量操作:支持批量插入和更新
  • 连接复用:高效的连接池管理

2. 预处理语句 (Prepared Statements)

这是 mysql2 最重要的安全特性:

code
┌──────────────┐        ┌──────────────┐
│  SQL 模板    │───────>│  MySQL 服务器 │
│  SELECT *    │        │  1. 解析      │
│  FROM users  │        │  2. 编译      │
│  WHERE id=?  │        │  3. 缓存      │
└──────────────┘        └──────────────┘
                              ▲
                              │
┌──────────────┐              │
│  参数值      │──────────────┘
│  [123]       │        执行(可重复)
└──────────────┘

优势:

  • 防止 SQL 注入:从根本上杜绝注入攻击
  • 提升性能:重复查询只需编译一次
  • 类型安全:自动处理参数类型转换

3. Promise 与 Async/Await

现代化的异步编程支持:

javascript
// 传统回调方式
pool.query('SELECT * FROM users', (err, results) => {
  // ...
})

// Promise 方式
const [results] = await pool.execute('SELECT * FROM users')

4. 流式查询

处理大数据集的利器:

code
大数据集查询 (100万行)
    │
    ├─ 传统方式:一次性加载到内存 → ❌ 内存溢出风险
    │
    └─ 流式查询:逐行处理 → ✅ 恒定内存占用

5. TypeScript 支持

开箱即用的类型定义,提供良好的开发体验。


安装与配置

安装

bash
# 安装 mysql2
npm install mysql2

# 可选:安装环境变量管理
npm install dotenv

配置参数详解

基础连接参数

javascript
const mysql = require('mysql2')

const connection = mysql.createConnection({
  // 必需参数
  host: 'localhost',           // 数据库主机地址
  user: 'root',                // 数据库用户名
  password: 'password',        // 数据库密码
  database: 'mydb',            // 数据库名称
  
  // 可选参数
  port: 3306,                  // 端口号,默认 3306
  charset: 'utf8mb4',          // 字符集,默认 utf8mb4
  timezone: '+08:00',          // 时区设置
  connectTimeout: 10000,       // 连接超时时间(毫秒)
  
  // SSL 配置
  ssl: {
    ca: fs.readFileSync('/path/to/ca.pem'),
    cert: fs.readFileSync('/path/to/client-cert.pem'),
    key: fs.readFileSync('/path/to/client-key.pem')
  }
})

连接池参数

javascript
const pool = mysql.createPool({
  // 基础参数
  host: 'localhost',
  user: 'root',
  password: 'password',
  database: 'mydb',
  
  // 连接池配置
  connectionLimit: 10,         // 最大连接数(默认 10)
  waitForConnections: true,    // 连接池满时是否等待(默认 true)
  queueLimit: 0,               // 等待队列最大长度(0 表示无限制)
  
  // 连接生命周期管理
  maxIdle: 10,                 // 最大空闲连接数
  idleTimeout: 60000,          // 空闲连接超时时间(毫秒)
  
  // 连接保活
  enableKeepAlive: true,       // 启用 TCP Keep-Alive
  keepAliveInitialDelay: 0,    // Keep-Alive 初始延迟
  
  // 性能优化
  multipleStatements: false,   // 禁用多语句(安全考虑)
  namedPlaceholders: false,    // 启用命名占位符
})

参数最佳实践:

参数推荐值说明
connectionLimit10-20根据 CPU 核心数和数据库配置调整
queueLimit0无限制等待,配合应用层超时控制
waitForConnectionstrue避免连接池满时立即失败
idleTimeout30000-60000及时释放空闲连接
multipleStatementsfalse安全考虑,禁用多语句执行

环境变量配置

创建 .env 文件:

env
# 数据库配置
DB_HOST=localhost
DB_PORT=3306
DB_USER=root
DB_PASSWORD=your_secure_password
DB_NAME=myapp_database

# 连接池配置
DB_CONNECTION_LIMIT=10
DB_QUEUE_LIMIT=0

# SSL 配置(生产环境)
DB_SSL_ENABLED=true
DB_SSL_CA=/path/to/ca.pem

加载配置:

javascript
require('dotenv').config()

const mysql = require('mysql2/promise')

const pool = mysql.createPool({
  host: process.env.DB_HOST || 'localhost',
  port: parseInt(process.env.DB_PORT) || 3306,
  user: process.env.DB_USER || 'root',
  password: process.env.DB_PASSWORD,
  database: process.env.DB_NAME,
  connectionLimit: parseInt(process.env.DB_CONNECTION_LIMIT) || 10,
  queueLimit: parseInt(process.env.DB_QUEUE_LIMIT) || 0
})

module.exports = pool

快速入门

基础连接示例

javascript
const mysql = require('mysql2')

// 创建连接
const connection = mysql.createConnection({
  host: 'localhost',
  user: 'root',
  password: 'password',
  database: 'testdb'
})

// 执行查询
connection.query('SELECT * FROM users LIMIT 5', (err, results, fields) => {
  if (err) {
    console.error('查询错误:', err.message)
    return
  }
  console.log('查询结果:', results)
  console.log('字段信息:', fields.map(f => f.name))
})

// 关闭连接
connection.end()

使用连接池(推荐)

javascript
const mysql = require('mysql2')
const pool = mysql.createPool({
  host: 'localhost',
  user: 'root',
  password: 'password',
  database: 'testdb',
  connectionLimit: 10
})

// 直接使用 pool.query(自动管理连接)
pool.query('SELECT * FROM products', (err, results) => {
  if (err) throw err
  console.log('产品列表:', results)
})

// 关闭连接池
pool.end(err => {
  if (err) console.error('关闭连接池失败:', err)
  console.log('连接池已关闭')
})

Promise 版本(推荐)

javascript
const mysql = require('mysql2/promise')

async function main() {
  // 创建连接池
  const pool = mysql.createPool({
    host: 'localhost',
    user: 'root',
    password: 'password',
    database: 'testdb'
  })

  try {
    // 执行查询
    const [rows, fields] = await pool.execute(
      'SELECT * FROM users WHERE age > ?',
      [18]
    )
    console.log('成年用户:', rows)
  } catch (error) {
    console.error('查询失败:', error)
  } finally {
    await pool.end()
  }
}

main()

核心概念

连接 vs 连接池

code
┌─────────────────────────────────────────────────────────────┐
│                    单连接模式                                │
│  ┌────────┐      ┌──────────┐      ┌──────────┐           │
│  │ 请求 1 │─────>│ 等待连接 │─────>│ 执行查询 │           │
│  └────────┘      └──────────┘      └──────────┘           │
│  ┌────────┐                        ┌──────────┐           │
│  │ 请求 2 │───────────────────────>│   阻塞   │ ❌        │
│  └────────┘                        └──────────┘           │
└─────────────────────────────────────────────────────────────┘

┌─────────────────────────────────────────────────────────────┐
│                    连接池模式                                │
│  ┌────────┐      ┌──────────────────────┐                  │
│  │ 请求 1 │─────>│ 从池中获取连接       │──┐               │
│  └────────┘      └──────────────────────┘  │               │
│  ┌────────┐      ┌──────────────────────┐  │  ┌─────────┐ │
│  │ 请求 2 │─────>│ 从池中获取另一个连接 │──┼─>│ 并发执行│ ✅│
│  └────────┘      └──────────────────────┘  │  └─────────┘ │
│  ┌────────┐      ┌──────────────────────┐  │               │
│  │ 请求 3 │─────>│ 从池中获取连接       │──┘               │
│  └────────┘      └──────────────────────┘                  │
└─────────────────────────────────────────────────────────────┘

对比表:

特性单连接连接池
并发能力❌ 低(串行执行)✅ 高(并行执行)
连接开销❌ 高(频繁创建销毁)✅ 低(连接复用)
资源管理❌ 手动管理✅ 自动管理
适用场景脚本、简单应用Web 服务、生产环境

query vs execute

方法工作原理安全性性能适用场景
query()客户端拼接 SQL 后发送⚠️ 需手动防注入中等流式查询、一次性查询
execute()先发送模板,再发送参数✅ 自动防注入高(可缓存)参数化查询、重复查询

示例对比:

javascript
// ❌ 危险:直接拼接 SQL(SQL 注入风险)
const userId = "1 OR 1=1"
pool.query(`SELECT * FROM users WHERE id = ${userId}`)

// ⚠️ 使用 query + 参数(安全但无预处理优势)
pool.query('SELECT * FROM users WHERE id = ?', [userId])

// ✅ 推荐:使用 execute(预处理语句)
pool.execute('SELECT * FROM users WHERE id = ?', [userId])

结果集结构

javascript
const [rows, fields] = await pool.execute('SELECT id, name FROM users')

// rows: 结果数据数组
[
  { id: 1, name: 'Alice' },
  { id: 2, name: 'Bob' }
]

// fields: 字段元数据数组
[
  {
    name: 'id',
    type: 'LONG',
    length: 11,
    db: 'mydb',
    table: 'users',
    ...
  },
  {
    name: 'name',
    type: 'VAR_STRING',
    length: 255,
    ...
  }
]

API 接口说明

连接池 API

pool.execute(sql, [values])

执行预处理语句(推荐)。

javascript
// 查询
const [rows] = await pool.execute('SELECT * FROM users WHERE id = ?', [1])

// 插入
const [result] = await pool.execute(
  'INSERT INTO users (name, email) VALUES (?, ?)',
  ['Alice', 'alice@example.com']
)
console.log('插入ID:', result.insertId)

// 更新
const [result] = await pool.execute(
  'UPDATE users SET name = ? WHERE id = ?',
  ['Bob', 1]
)
console.log('影响行数:', result.affectedRows)

// 删除
const [result] = await pool.execute('DELETE FROM users WHERE id = ?', [1])
console.log('删除行数:', result.affectedRows)

pool.query(sql, [values])

执行普通查询。

javascript
const [rows] = await pool.query('SELECT NOW() as now')

pool.getConnection

从连接池获取连接(用于事务)。

javascript
const connection = await pool.getConnection()
try {
  await connection.beginTransaction()
  // ... 执行多个操作
  await connection.commit()
} catch (error) {
  await connection.rollback()
  throw error
} finally {
  connection.release() // 归还连接
}

pool.end

关闭连接池。

javascript
await pool.end()

连接 API

connection.beginTransaction

开始事务。

connection.commit

提交事务。

connection.rollback

回滚事务。

connection.release

释放连接回连接池。

结果对象属性

INSERT 操作结果

javascript
{
  affectedRows: 1,      // 影响的行数
  insertId: 123,        // 插入的自增 ID
  warningStatus: 0      // 警告状态
}

UPDATE/DELETE 操作结果

javascript
{
  affectedRows: 2,      // 影响的行数
  changedRows: 2,       // 实际修改的行数
  warningStatus: 0
}

SELECT 操作结果

javascript
// 返回数组,每行一个对象
[
  { id: 1, name: 'Alice', email: 'alice@example.com' },
  { id: 2, name: 'Bob', email: 'bob@example.com' }
]

高级特性

1. 流式查询

处理大数据集时,避免内存溢出:

javascript
const fs = require('fs')
const pool = require('./database')

async function exportLargeDataset() {
  return new Promise((resolve, reject) => {
    const writeStream = fs.createWriteStream('export.csv')
    const queryStream = pool
      .query('SELECT id, name, email FROM users')
      .stream()

    writeStream.write('id,name,email\n')

    queryStream.on('data', (row) => {
      // 每行数据处理
      writeStream.write(`${row.id},"${row.name}","${row.email}"\n`)
    })

    queryStream.on('end', () => {
      writeStream.end()
      console.log('导出完成')
      resolve()
    })

    queryStream.on('error', (err) => {
      writeStream.end()
      reject(err)
    })
  })
}

2. 事务处理

完整的事务处理示例:

javascript
async function transferMoney(fromId, toId, amount) {
  const connection = await pool.getConnection()
  
  try {
    await connection.beginTransaction()

    // 检查余额
    const [accounts] = await connection.execute(
      'SELECT balance FROM accounts WHERE id = ? FOR UPDATE',
      [fromId]
    )
    
    if (accounts.length === 0) {
      throw new Error('源账户不存在')
    }
    
    if (accounts[0].balance < amount) {
      throw new Error('余额不足')
    }

    // 扣款
    await connection.execute(
      'UPDATE accounts SET balance = balance - ? WHERE id = ?',
      [amount, fromId]
    )

    // 加款
    await connection.execute(
      'UPDATE accounts SET balance = balance + ? WHERE id = ?',
      [amount, toId]
    )

    // 记录交易
    await connection.execute(
      'INSERT INTO transactions (from_id, to_id, amount) VALUES (?, ?, ?)',
      [fromId, toId, amount]
    )

    await connection.commit()
    return { success: true, message: '转账成功' }
    
  } catch (error) {
    await connection.rollback()
    return { success: false, error: error.message }
  } finally {
    connection.release()
  }
}

3. 批量操作

javascript
// 批量插入
async function batchInsertUsers(users) {
  const sql = 'INSERT INTO users (name, email, age) VALUES ?'
  const values = users.map(u => [u.name, u.email, u.age])
  
  const [result] = await pool.query(sql, [values])
  console.log(`成功插入 ${result.affectedRows} 条记录`)
  return result
}

// 使用示例
await batchInsertUsers([
  { name: 'Alice', email: 'alice@example.com', age: 25 },
  { name: 'Bob', email: 'bob@example.com', age: 30 },
  { name: 'Charlie', email: 'charlie@example.com', age: 28 }
])

4. TypeScript 支持

typescript
import { RowDataPacket, ResultSetHeader, Pool } from 'mysql2/promise'

// 定义数据类型
interface User extends RowDataPacket {
  id: number
  name: string
  email: string
  age: number
  createdAt: Date
}

// 类型安全的查询
async function getUserById(pool: Pool, id: number): Promise<User | null> {
  const [rows] = await pool.execute<User[]>(
    'SELECT * FROM users WHERE id = ?',
    [id]
  )
  return rows.length > 0 ? rows[0] : null
}

// 类型安全的插入
async function createUser(
  pool: Pool, 
  user: Omit<User, 'id' | 'createdAt'>
): Promise<number> {
  const [result] = await pool.execute<ResultSetHeader>(
    'INSERT INTO users (name, email, age) VALUES (?, ?, ?)',
    [user.name, user.email, user.age]
  )
  return result.insertId
}

5. 命名占位符

javascript
const mysql = require('mysql2/promise')

const pool = mysql.createPool({
  host: 'localhost',
  user: 'root',
  password: 'password',
  database: 'testdb',
  namedPlaceholders: true  // 启用命名占位符
})

// 使用命名参数
const [rows] = await pool.execute(
  'SELECT * FROM users WHERE age > :minAge AND city = :city',
  { minAge: 18, city: 'Beijing' }
)

性能优化

1. 连接池优化

优化策略:

javascript
// 根据 CPU 核心数动态调整
const os = require('os')
const cpuCount = os.cpus().length

const pool = mysql.createPool({
  host: 'localhost',
  user: 'root',
  password: 'password',
  database: 'testdb',
  
  // 连接数 = CPU 核心数 * 2 + 1
  connectionLimit: cpuCount * 2 + 1,
  
  // 空闲连接管理
  maxIdle: cpuCount,
  idleTimeout: 30000,
  
  // 等待队列
  waitForConnections: true,
  queueLimit: 100
})

监控连接池状态:

javascript
// 定期监控
setInterval(() => {
  console.log({
    totalConnections: pool.pool._allConnections.length,
    freeConnections: pool.pool._freeConnections.length,
    waitingForConnection: pool.pool._connectionQueue.length
  })
}, 5000)

2. 查询缓存

使用 Redis 缓存热点数据:

javascript
const redis = require('redis')
const redisClient = redis.createClient({ url: process.env.REDIS_URL })
await redisClient.connect()

async function getProductWithCache(productId) {
  const cacheKey = `product:${productId}`
  
  // 1. 检查缓存
  const cached = await redisClient.get(cacheKey)
  if (cached) {
    return JSON.parse(cached)
  }
  
  // 2. 查询数据库
  const [rows] = await pool.execute(
    'SELECT * FROM products WHERE id = ?',
    [productId]
  )
  const product = rows[0]
  
  // 3. 写入缓存
  if (product) {
    await redisClient.setEx(cacheKey, 300, JSON.stringify(product))
  }
  
  return product
}

3. 慢查询监控

javascript
// 慢查询监控中间件
async function executeWithMonitoring(sql, params, threshold = 500) {
  const start = Date.now()
  
  try {
    const result = await pool.execute(sql, params)
    return result
  } finally {
    const duration = Date.now() - start
    
    if (duration > threshold) {
      console.warn({
        type: 'SLOW_QUERY',
        duration: `${duration}ms`,
        sql,
        params,
        timestamp: new Date().toISOString()
      })
    }
  }
}

4. 批量操作优化

javascript
// 批量插入优化
async function batchInsertOptimized(items) {
  const BATCH_SIZE = 1000
  const results = []
  
  for (let i = 0; i < items.length; i += BATCH_SIZE) {
    const batch = items.slice(i, i + BATCH_SIZE)
    const values = batch.map(item => [item.name, item.email])
    
    const sql = 'INSERT INTO users (name, email) VALUES ?'
    const [result] = await pool.query(sql, [values])
    results.push(result)
  }
  
  return results
}

安全最佳实践

1. 防止 SQL 注入

javascript
// ❌ 错误:直接拼接 SQL
const userId = req.params.id
pool.query(`SELECT * FROM users WHERE id = ${userId}`)  // 危险!

// ✅ 正确:使用预处理语句
pool.execute('SELECT * FROM users WHERE id = ?', [userId])

// ✅ 正确:使用命名占位符
pool.execute('SELECT * FROM users WHERE id = :id', { id: userId })

2. 敏感信息管理

javascript
// ❌ 错误:硬编码密码
const pool = mysql.createPool({
  password: 'my_password_123'  // 危险!
})

// ✅ 正确:使用环境变量
const pool = mysql.createPool({
  password: process.env.DB_PASSWORD
})

3. 最小权限原则

sql
-- 创建应用专用用户
CREATE USER 'app_user'@'%' IDENTIFIED BY 'secure_password';

-- 只授予必要的权限
GRANT SELECT, INSERT, UPDATE ON myapp.users TO 'app_user'@'%';
GRANT SELECT ON myapp.products TO 'app_user'@'%';

-- 禁止危险操作
REVOKE DELETE ON myapp.* FROM 'app_user'@'%';

FLUSH PRIVILEGES;

4. SSL/TLS 加密

javascript
const fs = require('fs')

const pool = mysql.createPool({
  host: process.env.DB_HOST,
  user: process.env.DB_USER,
  password: process.env.DB_PASSWORD,
  database: process.env.DB_NAME,
  
  ssl: {
    ca: fs.readFileSync('/path/to/ca.pem'),
    cert: fs.readFileSync('/path/to/client-cert.pem'),
    key: fs.readFileSync('/path/to/client-key.pem'),
    rejectUnauthorized: true  // 严格验证证书
  }
})

5. 禁用多语句执行

javascript
// ❌ 危险:启用多语句
const pool = mysql.createPool({
  multipleStatements: true  // 允许执行多条 SQL,增加注入风险
})

// ✅ 安全:禁用多语句(默认值)
const pool = mysql.createPool({
  multipleStatements: false
})

错误处理与调试

常见错误代码

错误代码说明解决方案
ECONNREFUSED无法连接数据库检查数据库服务是否启动
ER_ACCESS_DENIED_ERROR访问被拒绝检查用户名密码
ER_DUP_ENTRY唯一键冲突检查数据是否重复
ER_NO_REFERENCED_ROW_2外键约束失败检查关联数据是否存在
ER_PARSE_ERRORSQL 语法错误检查 SQL 语法
PROTOCOL_CONNECTION_LOST连接丢失实现重连机制
ER_LOCK_DEADLOCK死锁实现重试机制

错误处理示例

javascript
class DatabaseError extends Error {
  constructor(message, code, sql = '') {
    super(message)
    this.name = 'DatabaseError'
    this.code = code
    this.sql = sql
  }
}

async function safeQuery(sql, params) {
  try {
    return await pool.execute(sql, params)
  } catch (error) {
    // 记录详细错误信息
    console.error({
      message: error.message,
      code: error.code,
      sql: error.sql,
      params,
      timestamp: new Date().toISOString()
    })
    
    // 转换为应用错误
    switch (error.code) {
      case 'ER_DUP_ENTRY':
        throw new DatabaseError('数据已存在', 'DUPLICATE_ENTRY', sql)
      case 'ER_NO_REFERENCED_ROW_2':
        throw new DatabaseError('关联数据不存在', 'INVALID_REFERENCE', sql)
      default:
        throw new DatabaseError('数据库操作失败', 'DATABASE_ERROR', sql)
    }
  }
}

自动重试机制

javascript
const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms))

async function withRetry(fn, options = {}) {
  const {
    maxRetries = 3,
    initialDelay = 100,
    shouldRetry = (err) => ['ER_LOCK_DEADLOCK', 'PROTOCOL_CONNECTION_LOST'].includes(err.code)
  } = options
  
  for (let attempt = 1; attempt <= maxRetries; attempt++) {
    try {
      return await fn()
    } catch (error) {
      if (attempt === maxRetries || !shouldRetry(error)) {
        throw error
      }
      
      const delay = initialDelay * Math.pow(2, attempt - 1)
      console.warn(`尝试 ${attempt}/${maxRetries} 失败,${delay}ms 后重试...`)
      await sleep(delay)
    }
  }
}

// 使用示例
const result = await withRetry(
  () => pool.execute('UPDATE accounts SET balance = balance - ? WHERE id = ?', [100, 1]),
  { maxRetries: 3 }
)

实战应用

Express.js REST API 完整示例

项目结构:

code
project/
├── src/
│   ├── config/
│   │   └── database.js
│   ├── models/
│   │   └── user.model.js
│   ├── controllers/
│   │   └── user.controller.js
│   ├── routes/
│   │   └── user.routes.js
│   ├── middleware/
│   │   └── errorHandler.js
│   └── app.js
├── .env
└── package.json

数据库配置 (src/config/database.js):

javascript
require('dotenv').config()
const mysql = require('mysql2/promise')

const pool = mysql.createPool({
  host: process.env.DB_HOST || 'localhost',
  port: parseInt(process.env.DB_PORT) || 3306,
  user: process.env.DB_USER || 'root',
  password: process.env.DB_PASSWORD,
  database: process.env.DB_NAME,
  waitForConnections: true,
  connectionLimit: 10,
  queueLimit: 0,
  enableKeepAlive: true,
  keepAliveInitialDelay: 0
})

// 监听连接池事件
pool.on('acquire', (connection) => {
  console.log(`连接 ${connection.threadId} 被获取`)
})

pool.on('release', (connection) => {
  console.log(`连接 ${connection.threadId} 被释放`)
})

pool.on('enqueue', () => {
  console.warn('连接池已满,请求正在排队')
})

module.exports = pool

用户模型 (src/models/user.model.js):

javascript
const pool = require('../config/database')

class UserModel {
  async findAll(options = {}) {
    const { page = 1, limit = 10 } = options
    const offset = (page - 1) * limit
    
    const [rows] = await pool.execute(
      'SELECT id, name, email, created_at FROM users ORDER BY created_at DESC LIMIT ? OFFSET ?',
      [limit, offset]
    )
    
    const [countRows] = await pool.execute('SELECT COUNT(*) as total FROM users')
    
    return {
      data: rows,
      pagination: {
        page,
        limit,
        total: countRows[0].total,
        totalPages: Math.ceil(countRows[0].total / limit)
      }
    }
  }

  async findById(id) {
    const [rows] = await pool.execute(
      'SELECT id, name, email, created_at FROM users WHERE id = ?',
      [id]
    )
    return rows[0] || null
  }

  async create(userData) {
    const { name, email, password } = userData
    
    const [result] = await pool.execute(
      'INSERT INTO users (name, email, password) VALUES (?, ?, ?)',
      [name, email, password]
    )
    
    return this.findById(result.insertId)
  }

  async update(id, userData) {
    const fields = []
    const values = []
    
    if (userData.name) {
      fields.push('name = ?')
      values.push(userData.name)
    }
    if (userData.email) {
      fields.push('email = ?')
      values.push(userData.email)
    }
    
    if (fields.length === 0) {
      return this.findById(id)
    }
    
    values.push(id)
    
    await pool.execute(
      `UPDATE users SET ${fields.join(', ')} WHERE id = ?`,
      values
    )
    
    return this.findById(id)
  }

  async delete(id) {
    const [result] = await pool.execute('DELETE FROM users WHERE id = ?', [id])
    return result.affectedRows > 0
  }
}

module.exports = new UserModel()

用户控制器 (src/controllers/user.controller.js):

javascript
const userModel = require('../models/user.model')

class UserController {
  async getUsers(req, res, next) {
    try {
      const result = await userModel.findAll(req.query)
      res.json(result)
    } catch (error) {
      next(error)
    }
  }

  async getUser(req, res, next) {
    try {
      const user = await userModel.findById(req.params.id)
      if (!user) {
        return res.status(404).json({ error: '用户不存在' })
      }
      res.json(user)
    } catch (error) {
      next(error)
    }
  }

  async createUser(req, res, next) {
    try {
      const user = await userModel.create(req.body)
      res.status(201).json(user)
    } catch (error) {
      if (error.code === 'ER_DUP_ENTRY') {
        return res.status(409).json({ error: '邮箱已被注册' })
      }
      next(error)
    }
  }

  async updateUser(req, res, next) {
    try {
      const user = await userModel.update(req.params.id, req.body)
      if (!user) {
        return res.status(404).json({ error: '用户不存在' })
      }
      res.json(user)
    } catch (error) {
      next(error)
    }
  }

  async deleteUser(req, res, next) {
    try {
      const deleted = await userModel.delete(req.params.id)
      if (!deleted) {
        return res.status(404).json({ error: '用户不存在' })
      }
      res.status(204).send()
    } catch (error) {
      next(error)
    }
  }
}

module.exports = new UserController()

路由定义 (src/routes/user.routes.js):

javascript
const express = require('express')
const userController = require('../controllers/user.controller')

const router = express.Router()

router.get('/', userController.getUsers)
router.get('/:id', userController.getUser)
router.post('/', userController.createUser)
router.put('/:id', userController.updateUser)
router.delete('/:id', userController.deleteUser)

module.exports = router

错误处理中间件 (src/middleware/errorHandler.js):

javascript
function errorHandler(err, req, res, next) {
  console.error('错误:', err)

  if (err.code === 'ER_DUP_ENTRY') {
    return res.status(409).json({ error: '数据重复' })
  }

  if (err.code === 'ER_NO_REFERENCED_ROW_2') {
    return res.status(400).json({ error: '关联数据不存在' })
  }

  if (err.name === 'DatabaseError') {
    return res.status(500).json({ error: '数据库操作失败' })
  }

  res.status(500).json({ error: '服务器内部错误' })
}

module.exports = errorHandler

应用入口 (src/app.js):

javascript
require('dotenv').config()
const express = require('express')
const userRoutes = require('./routes/user.routes')
const errorHandler = require('./middleware/errorHandler')

const app = express()

// 中间件
app.use(express.json())

// 路由
app.use('/api/users', userRoutes)

// 错误处理
app.use(errorHandler)

// 启动服务器
const PORT = process.env.PORT || 3000
app.listen(PORT, () => {
  console.log(`服务器运行在端口 ${PORT}`)
})

常见问题解答

Q1: 连接池应该设置多少个连接?

A: 推荐公式:连接数 = CPU 核心数 × 2 + 1

javascript
const os = require('os')
const connectionLimit = os.cpus().length * 2 + 1

实际应用中需要通过压力测试调整最优值。

Q2: query 和 execute 该用哪个?

A: 推荐使用 execute

  • ✅ 自动防止 SQL 注入
  • ✅ 重复查询性能更好(预处理语句缓存)
  • ✅ 类型安全

只有在需要流式查询时才使用 query

Q3: 如何处理数据库连接丢失?

A: 实现自动重连和重试机制:

javascript
const pool = mysql.createPool({
  // ... 其他配置
  enableKeepAlive: true,
  keepAliveInitialDelay: 30000  // 30秒发送一次心跳
})

Q4: 如何调试 SQL 查询?

A: 启用查询日志:

javascript
// 开发环境启用查询日志
if (process.env.NODE_ENV === 'development') {
  const originalExecute = pool.execute.bind(pool)
  pool.execute = async function(sql, params) {
    console.log('SQL:', sql)
    console.log('参数:', params)
    const start = Date.now()
    const result = await originalExecute(sql, params)
    console.log(`执行时间: ${Date.now() - start}ms`)
    return result
  }
}

Q5: 如何处理大结果集?

A: 使用流式查询:

javascript
const stream = pool.query('SELECT * FROM large_table').stream()

stream.on('data', (row) => {
  // 逐行处理
})

stream.on('end', () => {
  console.log('处理完成')
})

Q6: 如何实现数据库迁移?

A: 建议使用专门的迁移工具:

Q7: 如何测试数据库代码?

A: 使用测试数据库和事务回滚:

javascript
describe('UserService', () => {
  let connection
  
  beforeAll(async () => {
    connection = await pool.getConnection()
  })
  
  afterAll(async () => {
    connection.release()
  })
  
  beforeEach(async () => {
    await connection.beginTransaction()
  })
  
  afterEach(async () => {
    await connection.rollback()
  })
  
  test('创建用户', async () => {
    const user = await createUser({ name: 'Test' })
    expect(user.name).toBe('Test')
  })
})

选择建议

mysql2 vs ORM

特性mysql2 (原生驱动)ORM (如 Sequelize)
抽象层次低。直接编写 SQL高。面向对象操作
控制力极高。完全控制 SQL中等。自动生成 SQL
开发速度较慢。手动编写 SQL快。内置 CRUD
学习曲线低。熟悉 SQL 即可较高。需学习 ORM API
性能通常更高。无中间层可能较低。SQL 不是最优
类型安全需手动定义自动生成类型

适用场景

选择 mysql2 当:

  • ✅ 追求极致性能
  • ✅ 需要完全控制 SQL
  • ✅ 项目较简单
  • ✅ 团队熟悉 SQL

选择 ORM 当:

  • ✅ 复杂的业务模型
  • ✅ 需要快速开发
  • ✅ 需要数据库迁移功能
  • ✅ 希望代码与数据库解耦

参考资料